ποΈGitΠ―ΡΠ°ποΈ
Node / meshtastic / Meshtastic-Android / files / core / takserver / src / commonMain / kotlin / org / meshtastic / core / takserver / TAKMeshIntegration.kt
Displaying Raw β’ Download
core/takserver/src/commonMain/kotlin/org/meshtastic/core/takserver/TAKMeshIntegration.kt b4bedd92fcc8d65daff79200d71f6ebc3d8f0cfa (b4bedd92) Text, 25.05 KB
T8b949e/*
* Copyright (c) 2026 Meshtastic LLC
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program. If not, see <https://www.gnu.org/licenses/>.
*/
Tf0883e@fileTb4b4b4:Te6edf3SuppressTb4b4b4(Ta5d6ff"Ta5d6ffReturnCountTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffLongMethodTa5d6ff"Tb4b4b4)
Tff7b72package T7ee787org.meshtastic.core.takserver
Tff7b72import T7ee787co.touchlab.kermit.Logger
Tff7b72import T7ee787kotlinx.coroutines.CoroutineScope
Tff7b72import T7ee787kotlinx.coroutines.Job
Tff7b72import T7ee787kotlinx.coroutines.flow.distinctUntilChanged
Tff7b72import T7ee787kotlinx.coroutines.flow.filter
Tff7b72import T7ee787kotlinx.coroutines.flow.map
Tff7b72import T7ee787kotlinx.coroutines.launch
Tff7b72import T7ee787okio.ByteString.Companion.toByteString
Tff7b72import T7ee787org.meshtastic.core.model.Capabilities
Tff7b72import T7ee787org.meshtastic.core.model.DataPacket
Tff7b72import T7ee787org.meshtastic.core.model.NodeAddress
Tff7b72import T7ee787org.meshtastic.core.repository.CommandSender
Tff7b72import T7ee787org.meshtastic.core.repository.MeshConfigHandler
Tff7b72import T7ee787org.meshtastic.core.repository.NodeRepository
Tff7b72import T7ee787org.meshtastic.core.repository.ServiceRepository
Tff7b72import T7ee787org.meshtastic.core.takserver.TAKPacketConversion.toTAKPacket
Tff7b72import T7ee787org.meshtastic.core.takserver.TAKPacketV2Conversion.toTAKPacketV2
Tff7b72import T7ee787org.meshtastic.proto.MemberRole
Tff7b72import T7ee787org.meshtastic.proto.MeshPacket
Tff7b72import T7ee787org.meshtastic.proto.PortNum
Tff7b72import T7ee787org.meshtastic.proto.TAKPacket
Tff7b72import T7ee787org.meshtastic.proto.Team
Tff7b72import T7ee787kotlin.concurrent.Volatile
Tff7b72import T7ee787kotlin.concurrent.atomics.AtomicBoolean
Tff7b72import T7ee787kotlin.concurrent.atomics.ExperimentalAtomicApi
Tff7b72import T7ee787kotlin.random.Random
Tff7b72import T7ee787kotlin.time.Clock
Tff7b72import T7ee787kotlin.time.Duration.Companion.minutes
T8b949e/**
* Bidirectional bridge between the local TAK server and the Meshtastic mesh network.
*
* Outbound traffic (TAK client -> mesh) is version-gated on the connected radio's firmware version, exposed via
* [Capabilities.supportsTakV2]:
* - Firmware **>= 2.8.0**: TAKPacketV2 on port 78 (ATAK_PLUGIN_V2) with zstd dictionary compression via TAKPacket-SDK.
* Supports all CoT payload types (PLI, GeoChat, DrawnShape, Marker, Route, Aircraft, Casevac, Emergency, Task) with
* compact typed encodings that fit under the 237B LoRa MTU.
* - Firmware **<= 2.7.x**: Legacy [TAKPacket] on port 72 (ATAK_PLUGIN) with bare protobuf encoding. Supports only PLI
* and GeoChat β shapes, markers, routes, and other typed CoT events are dropped (with a warning) because the legacy
* schema cannot represent them.
*
* Inbound traffic (mesh -> TAK client) is always dual-path tolerant β both port 72 and port 78 are dispatched
* regardless of the local radio's firmware version, so a v2-capable node can still relay legacy v1 packets received
* from older nodes in mixed-firmware mesh deployments.
*/
Tf0883e@OptInTb4b4b4(Te6edf3ExperimentalAtomicApiTff7b72::Te6edf3classTb4b4b4)
Tff7b72class T56d364TAKMeshIntegrationTb4b4b4(
Tff7b72private Tff7b72val Te6edf3takServerManagerTb4b4b4: Te6edf3TAKServerManagerTb4b4b4,
Tff7b72private Tff7b72val Te6edf3commandSenderTb4b4b4: Te6edf3CommandSenderTb4b4b4,
Tff7b72private Tff7b72val Te6edf3serviceRepositoryTb4b4b4: Te6edf3ServiceRepositoryTb4b4b4,
Tff7b72private Tff7b72val Te6edf3meshConfigHandlerTb4b4b4: Te6edf3MeshConfigHandlerTb4b4b4,
Tff7b72private Tff7b72val Te6edf3nodeRepositoryTb4b4b4: Te6edf3NodeRepositoryTb4b4b4,
Tff7b72private Tff7b72val Te6edf3meshToCotBroadcasterTb4b4b4: Te6edf3MeshToCotBroadcasterTb4b4b4,
Tb4b4b4) Tb4b4b4{
Tff7b72private Tff7b72val Te6edf3isRunning Tff7b72= Te6edf3AtomicBooleanTb4b4b4(Tff7b72falseTb4b4b4)
T8b949e// Immutable list reference replaced atomically in start()/stop(); never mutated in-place.
T8b949e// @Volatile only guarantees visibility of the reference itself β any in-place mutation
T8b949e// would bypass the visibility guarantee and must not be added.
Tf0883e@Volatile Tff7b72private Tff7b72var Te6edf3jobsTb4b4b4: Te6edf3ListTff7b72<Te6edf3JobTff7b72> Tff7b72= Te6edf3emptyListTb4b4b4(Tb4b4b4)
Tf0883e@Volatile Tff7b72private Tff7b72var Te6edf3currentTeamTb4b4b4: Te6edf3Team Tff7b72= Te6edf3TeamTb4b4b4.Te6edf3Unspecifed_Color
Tf0883e@Volatile Tff7b72private Tff7b72var Te6edf3currentRoleTb4b4b4: Te6edf3MemberRole Tff7b72= Te6edf3MemberRoleTb4b4b4.Te6edf3Unspecifed
T8b949e// Drops CoT the bridge has already injected within a short window β guards against the
T8b949e// same mesh message arriving via multiple LoRa relay paths or a retransmit. Touched only
T8b949e// from the single meshPacketFlow collector coroutine (handleMeshPacket), so no locking.
Tff7b72private Tff7b72val Te6edf3deliveryDedup Tff7b72= Te6edf3CotDeliveryDedupTb4b4b4(Tb4b4b4)
Tff7b72fun Td2a8ffstartTb4b4b4(Te6edf3scopeTb4b4b4: Te6edf3CoroutineScopeTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3isRunningTb4b4b4.Te6edf3compareAndSetTb4b4b4(Te6edf3expectedValue Tff7b72= Tff7b72falseTb4b4b4, Te6edf3newValue Tff7b72= Tff7b72trueTb4b4b4)Tb4b4b4) Tff7b72return
Te6edf3takServerManagerTb4b4b4.Te6edf3startTb4b4b4(Te6edf3scopeTb4b4b4)
Tff7b72val Te6edf3newJobs Tff7b72=
Te6edf3listOfTb4b4b4(
T8b949e// Forward incoming CoT from TAK clients to mesh
Te6edf3scopeTb4b4b4.Te6edf3launch Tb4b4b4{
Te6edf3takServerManagerTb4b4b4.Te6edf3inboundMessagesTb4b4b4.Te6edf3collect Tb4b4b4{ Tb4b4b4(Te6edf3cotMessageTb4b4b4, Te6edf3clientInfoTb4b4b4) Tff7b72-Tff7b72>
T8b949e// Enrich GeoChat messages with the originating TAK client's
T8b949e// callsign when the message itself lacks one. This only applies
T8b949e// to messages FROM the connected TAK client β mesh-originated
T8b949e// messages flow through handleMeshPacket() instead.
Tff7b72val Te6edf3enriched Tff7b72=
Tff7b72if Tb4b4b4(
Te6edf3cotMessageTb4b4b4.Te6edf3type Tff7b72=Tff7b72= Ta5d6ff"Ta5d6ffb-t-fTa5d6ff" Tff7b72&Tff7b72&
Te6edf3cotMessageTb4b4b4.Te6edf3contactTff7b72?.Te6edf3callsignTb4b4b4.Te6edf3isNullOrEmptyTb4b4b4(Tb4b4b4) Tff7b72&Tff7b72&
Te6edf3clientInfoTff7b72?.Te6edf3callsign Tff7b72!Tff7b72= Tff7b72null
Tb4b4b4) Tb4b4b4{
Te6edf3cotMessageTb4b4b4.Te6edf3copyTb4b4b4(
Te6edf3contact Tff7b72=
Tb4b4b4(Te6edf3cotMessageTb4b4b4.Te6edf3contact Tff7b72?: Te6edf3CoTContactTb4b4b4(Te6edf3callsign Tff7b72= Ta5d6ff"Ta5d6ff"Tb4b4b4)Tb4b4b4)Tb4b4b4.Te6edf3copyTb4b4b4(
Te6edf3callsign Tff7b72= Te6edf3clientInfoTb4b4b4.Te6edf3callsignTb4b4b4,
Tb4b4b4)Tb4b4b4,
Tb4b4b4)
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3cotMessage
Tb4b4b4}
Te6edf3sendCoTToMeshTb4b4b4(Te6edf3enrichedTb4b4b4)
Tb4b4b4}
Tb4b4b4}Tb4b4b4,
T8b949e// Forward incoming ATAK packets from mesh to TAK clients
Te6edf3scopeTb4b4b4.Te6edf3launch Tb4b4b4{
Te6edf3serviceRepositoryTb4b4b4.Te6edf3meshPacketFlow
Tb4b4b4.Te6edf3filter Tb4b4b4{
Tffa657itTb4b4b4.Te6edf3decodedTff7b72?.Te6edf3portnum Tff7b72=Tff7b72= Te6edf3PortNumTb4b4b4.Te6edf3ATAK_PLUGIN_V2 Tff7b72|Tff7b72| Tffa657itTb4b4b4.Te6edf3decodedTff7b72?.Te6edf3portnum Tff7b72=Tff7b72= Te6edf3PortNumTb4b4b4.Te6edf3ATAK_PLUGIN
Tb4b4b4}
Tb4b4b4.Te6edf3collect Tb4b4b4{ Te6edf3packet Tff7b72-Tff7b72> Te6edf3handleMeshPacketTb4b4b4(Te6edf3packetTb4b4b4) Tb4b4b4}
Tb4b4b4}Tb4b4b4,
T8b949e// Track TAK config changes
Te6edf3scopeTb4b4b4.Te6edf3launch Tb4b4b4{
Te6edf3meshConfigHandlerTb4b4b4.Te6edf3moduleConfig
Tb4b4b4.Te6edf3map Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3tak Tb4b4b4}
Tb4b4b4.Te6edf3distinctUntilChangedTb4b4b4(Tb4b4b4)
Tb4b4b4.Te6edf3collect Tb4b4b4{ Te6edf3takConfig Tff7b72-Tff7b72>
Te6edf3currentTeam Tff7b72= Te6edf3takConfigTff7b72?.Te6edf3team Tff7b72?: Te6edf3TeamTb4b4b4.Te6edf3Unspecifed_Color
Te6edf3currentRole Tff7b72= Te6edf3takConfigTff7b72?.Te6edf3role Tff7b72?: Te6edf3MemberRoleTb4b4b4.Te6edf3Unspecifed
Tb4b4b4}
Tb4b4b4}Tb4b4b4,
Tb4b4b4)
Te6edf3jobs Tff7b72= Te6edf3newJobs
T8b949e// Node -> CoT contacts. Self-gates on its own opt-in pref; starting it here means it can
T8b949e// only ever run while the TAK server is enabled.
Te6edf3meshToCotBroadcasterTb4b4b4.Te6edf3startTb4b4b4(Te6edf3scopeTb4b4b4)
Tff7b72val Te6edf3fw Tff7b72= Te6edf3nodeRepositoryTb4b4b4.Te6edf3myNodeInfoTb4b4b4.Te6edf3valueTff7b72?.Te6edf3firmwareVersion
Tff7b72val Te6edf3proto Tff7b72= Tff7b72if Tb4b4b4(Te6edf3CapabilitiesTb4b4b4(Te6edf3fwTb4b4b4)Tb4b4b4.Te6edf3supportsTakV2Tb4b4b4) Ta5d6ff"Ta5d6ffv2 (port 78, zstd)Ta5d6ff" Tff7b72else Ta5d6ff"Ta5d6ffv1 (port 72, legacy)Ta5d6ff"
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffTAK Mesh Integration started β firmware=Tffd700$Te6edf3fwTa5d6ff, outbound=Tffd700$Te6edf3protoTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tff7b72fun Td2a8ffstopTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3isRunningTb4b4b4.Te6edf3compareAndSetTb4b4b4(Te6edf3expectedValue Tff7b72= Tff7b72trueTb4b4b4, Te6edf3newValue Tff7b72= Tff7b72falseTb4b4b4)Tb4b4b4) Tff7b72return
Tff7b72val Te6edf3toCancel Tff7b72= Te6edf3jobs
Te6edf3jobs Tff7b72= Te6edf3emptyListTb4b4b4(Tb4b4b4)
Te6edf3toCancelTb4b4b4.Te6edf3forEachTb4b4b4(Te6edf3JobTff7b72::Te6edf3cancelTb4b4b4)
Te6edf3meshToCotBroadcasterTb4b4b4.Te6edf3stopTb4b4b4(Tb4b4b4)
Te6edf3takServerManagerTb4b4b4.Te6edf3stopTb4b4b4(Tb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffTAK Mesh Integration stoppedTa5d6ff" Tb4b4b4}
Tb4b4b4}
T8b949e// ββ Send: TAK client β mesh βββββββββββββββββββββββββββββββββββββββββββββ
T8b949e/**
* Determine the outbound TAK protocol version based on the connected radio's firmware version. Evaluated per-send
* (not cached) so the bridge picks up firmware upgrades during a session without restart. Until the version is
* known, use the legacy format which every TAK-capable firmware version supports.
*/
Tff7b72private Tff7b72fun Td2a8ffuseTakV2Tb4b4b4(Tb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
Tff7b72val Te6edf3fw Tff7b72= Te6edf3nodeRepositoryTb4b4b4.Te6edf3myNodeInfoTb4b4b4.Te6edf3valueTff7b72?.Te6edf3firmwareVersion Tff7b72?: Tff7b72return Tff7b72false
Tff7b72return Te6edf3CapabilitiesTb4b4b4(Te6edf3fwTb4b4b4)Tb4b4b4.Te6edf3supportsTakV2
Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffsendCoTToMeshTb4b4b4(Te6edf3cotMessageTb4b4b4: Te6edf3CoTMessageTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3useTakV2Tb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3sendCoTToMeshV2Tb4b4b4(Te6edf3cotMessageTb4b4b4)
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3sendCoTToMeshV1Tb4b4b4(Te6edf3cotMessageTb4b4b4)
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* v2 send path (firmware >= 2.8.0): SDK parser + zstd dictionary compression, full typed payload support
* (DrawnShape, Marker, Route, Aircraft, Casevac, Emergency, Task, plus PLI / GeoChat). Wire format: `[flags
* byte][zstd-compressed TAKPacketV2 protobuf]` on port 78 (ATAK_PLUGIN_V2).
*/
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffsendCoTToMeshV2Tb4b4b4(Te6edf3cotMessageTb4b4b4: Te6edf3CoTMessageTb4b4b4) Tb4b4b4{
T8b949e// Prefer the sourceEventXml for shape/marker/route types β the SDK's
T8b949e// CotXmlParser extracts compact typed payloads (DrawnShape, Marker,
T8b949e// Route, etc.) that compress far better than raw_detail encoding.
T8b949e// For PLI and GeoChat, use the enriched CoTMessage (which may have
T8b949e// had callsign/contact injected by the upstream enrichment step).
Tff7b72val Te6edf3rawXml Tff7b72= Te6edf3cotMessageTb4b4b4.Te6edf3sourceEventXml Tff7b72?: Te6edf3cotMessageTb4b4b4.Te6edf3toXmlTb4b4b4(Tb4b4b4)
T8b949e// Extend stale for static objects (routes, shapes, markers) that may
T8b949e// arrive over LoRa mesh past their original TTL. iTAK uses 2-min stale
T8b949e// for routes; ATAK uses 24h. 5 min ensures it survives mesh delivery.
Tff7b72val Te6edf3freshXml Tff7b72= Te6edf3ensureMinimumStaleForMeshTb4b4b4(Te6edf3rawXmlTb4b4b4)
T8b949e// Strip non-essential elements before compression to save wire bytes
Tff7b72val Te6edf3xml Tff7b72= Te6edf3stripNonEssentialElementsTb4b4b4(Te6edf3freshXmlTb4b4b4)
T8b949e// Route through the SDK parser/compressor which handles all typed
T8b949e// payloads (DrawnShape, Marker, Route, Aircraft, etc.) with compact
T8b949e// proto fields instead of raw_detail XML. Falls back to the app's
T8b949e// own conversion only if the SDK path fails.
T8b949e//
T8b949e// compressWithRemarksFallback preserves <remarks> text when the
T8b949e// compressed packet fits under the LoRa MTU, and strips remarks
T8b949e// automatically if needed to fit. Returns null if even without
T8b949e// remarks the packet exceeds the limit.
Tff7b72val Te6edf3wirePayloadTb4b4b4: Te6edf3ByteArray Tff7b72=
Tff7b72try Tb4b4b4{
Te6edf3TakSdkCompressorTb4b4b4.Te6edf3compressCoTTb4b4b4(Te6edf3xmlTb4b4b4, Te6edf3MAX_TAK_WIRE_PAYLOAD_BYTESTb4b4b4)
Tff7b72?: Te6edf3run Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{
Te6edf3buildString Tb4b4b4{
Te6edf3appendTb4b4b4(
Ta5d6ff"Ta5d6ffDropping oversized TAK packet: Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6fftype=Tffd700${Te6edf3cotMessageTb4b4b4.Te6edf3typeTffd700}Ta5d6ff max=Tffd700$Te6edf3MAX_TAK_WIRE_PAYLOAD_BYTESTa5d6ff"Tb4b4b4,
Tb4b4b4)
Te6edf3cotMessageTb4b4b4.Te6edf3sourceEventXmlTff7b72?.Te6edf3let Tb4b4b4{ Te6edf3src Tff7b72-Tff7b72>
Te6edf3appendTb4b4b4(Ta5d6ff'\n'Tb4b4b4)
Te6edf3appendTb4b4b4(Ta5d6ff"Ta5d6ffSource CoT event: Ta5d6ff"Tb4b4b4)
Te6edf3appendTb4b4b4(
Tff7b72if Tb4b4b4(Te6edf3srcTb4b4b4.Te6edf3length Tff7b72<Tff7b72= Te6edf3TAK_LOG_XML_MAX_CHARSTb4b4b4) Tb4b4b4{
Te6edf3src
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3srcTb4b4b4.Te6edf3takeTb4b4b4(Te6edf3TAK_LOG_XML_MAX_CHARSTb4b4b4) Tff7b72+ Ta5d6ff"Ta5d6ffβ¦Ta5d6ff"
Tb4b4b4}Tb4b4b4,
Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72return
Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3eTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffSDK parser/compressor failed for Tffd700${Te6edf3cotMessageTb4b4b4.Te6edf3typeTffd700}Ta5d6ff, trying app conversionTa5d6ff" Tb4b4b4}
Tff7b72val Te6edf3takPacketV2 Tff7b72= Te6edf3cotMessageTb4b4b4.Te6edf3toTAKPacketV2Tb4b4b4(Tb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3takPacketV2 Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffCannot convert CoT type Tffd700${Te6edf3cotMessageTb4b4b4.Te6edf3typeTffd700}Ta5d6ff to TAKPacketV2, droppingTa5d6ff" Tb4b4b4}
Tff7b72return
Tb4b4b4}
Tff7b72try Tb4b4b4{
Te6edf3TakV2CompressorTb4b4b4.Te6edf3compressTb4b4b4(Te6edf3takPacketV2Tb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3e2Tb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3e2Tb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffV2 compression failed for Tffd700${Te6edf3cotMessageTb4b4b4.Te6edf3typeTffd700}Ta5d6ff, using uncompressed wire formatTa5d6ff" Tb4b4b4}
Te6edf3encodeUncompressedTb4b4b4(Te6edf3takPacketV2Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tff7b72try Tb4b4b4{
Tff7b72val Te6edf3dataPacket Tff7b72=
Te6edf3DataPacketTb4b4b4(
Te6edf3to Tff7b72= Te6edf3NodeAddressTb4b4b4.Te6edf3ID_BROADCASTTb4b4b4,
Te6edf3bytes Tff7b72= Te6edf3wirePayloadTb4b4b4.Te6edf3toByteStringTb4b4b4(Tb4b4b4)Tb4b4b4,
Te6edf3dataType Tff7b72= Te6edf3PortNumTb4b4b4.Te6edf3ATAK_PLUGIN_V2Tb4b4b4.Te6edf3valueTb4b4b4,
Tb4b4b4)
Te6edf3commandSenderTb4b4b4.Te6edf3sendDataTb4b4b4(Te6edf3dataPacketTb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffSent V2 to mesh: Tffd700${Te6edf3cotMessageTb4b4b4.Te6edf3typeTffd700}Ta5d6ff (Tffd700${Te6edf3wirePayloadTb4b4b4.Te6edf3sizeTffd700}Ta5d6ff bytes)Ta5d6ff" Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3kotlinTb4b4b4.Te6edf3coroutinesTb4b4b4.Te6edf3cancellationTb4b4b4.Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3e
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
T8b949e// Something other than size β radio not connected, queue full, etc.
Te6edf3LoggerTb4b4b4.Te6edf3eTb4b4b4(Te6edf3eTb4b4b4) Tb4b4b4{
Ta5d6ff"Ta5d6ffFailed to send TAKPacketV2 to mesh (Tffd700${Te6edf3cotMessageTb4b4b4.Te6edf3typeTffd700}Ta5d6ff, Tffd700${Te6edf3wirePayloadTb4b4b4.Te6edf3sizeTffd700}Ta5d6ff bytes): Tffd700${Te6edf3eTb4b4b4.Te6edf3messageTffd700}Ta5d6ff"
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Legacy v1 send path (firmware <= 2.7.x): bare protobuf-encoded [TAKPacket] on port 72 (ATAK_PLUGIN), no zstd
* compression. Only PLI and GeoChat payloads are supported by the v1 schema β shapes, markers, routes, casevac,
* emergency, and task CoT events are dropped with a warning.
*/
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffsendCoTToMeshV1Tb4b4b4(Te6edf3cotMessageTb4b4b4: Te6edf3CoTMessageTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3takPacket Tff7b72=
Te6edf3cotMessageTb4b4b4.Te6edf3toTAKPacketTb4b4b4(Tb4b4b4)
Tff7b72?: Te6edf3run Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{
Ta5d6ff"Ta5d6ffDropping CoT for legacy v1 radio: type=Tffd700${Te6edf3cotMessageTb4b4b4.Te6edf3typeTffd700}Ta5d6ff not representable Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6ffin v1 TAKPacket schema (only PLI and GeoChat are supported). Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6ffUpgrade radio firmware to >= 2.8.0 for full payload support.Ta5d6ff"
Tb4b4b4}
Tff7b72return
Tb4b4b4}
Tff7b72val Te6edf3wirePayload Tff7b72= Te6edf3TAKPacketTb4b4b4.Te6edf3ADAPTERTb4b4b4.Te6edf3encodeTb4b4b4(Te6edf3takPacketTb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3wirePayloadTb4b4b4.Te6edf3size Tff7b72> Te6edf3MAX_TAK_WIRE_PAYLOAD_BYTESTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{
Ta5d6ff"Ta5d6ffDropping oversized v1 TAK packet: type=Tffd700${Te6edf3cotMessageTb4b4b4.Te6edf3typeTffd700}Ta5d6ff Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6ffsize=Tffd700${Te6edf3wirePayloadTb4b4b4.Te6edf3sizeTffd700}Ta5d6ffB max=Tffd700$Te6edf3MAX_TAK_WIRE_PAYLOAD_BYTESTa5d6ff"
Tb4b4b4}
Tff7b72return
Tb4b4b4}
Tff7b72try Tb4b4b4{
Tff7b72val Te6edf3dataPacket Tff7b72=
Te6edf3DataPacketTb4b4b4(
Te6edf3to Tff7b72= Te6edf3NodeAddressTb4b4b4.Te6edf3ID_BROADCASTTb4b4b4,
Te6edf3bytes Tff7b72= Te6edf3wirePayloadTb4b4b4.Te6edf3toByteStringTb4b4b4(Tb4b4b4)Tb4b4b4,
Te6edf3dataType Tff7b72= Te6edf3PortNumTb4b4b4.Te6edf3ATAK_PLUGINTb4b4b4.Te6edf3valueTb4b4b4,
Tb4b4b4)
Te6edf3commandSenderTb4b4b4.Te6edf3sendDataTb4b4b4(Te6edf3dataPacketTb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffSent V1 to mesh: Tffd700${Te6edf3cotMessageTb4b4b4.Te6edf3typeTffd700}Ta5d6ff (Tffd700${Te6edf3wirePayloadTb4b4b4.Te6edf3sizeTffd700}Ta5d6ff bytes)Ta5d6ff" Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3kotlinTb4b4b4.Te6edf3coroutinesTb4b4b4.Te6edf3cancellationTb4b4b4.Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3e
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3eTb4b4b4(Te6edf3eTb4b4b4) Tb4b4b4{
Ta5d6ff"Ta5d6ffFailed to send v1 TAKPacket to mesh (Tffd700${Te6edf3cotMessageTb4b4b4.Te6edf3typeTffd700}Ta5d6ff, Tffd700${Te6edf3wirePayloadTb4b4b4.Te6edf3sizeTffd700}Ta5d6ff bytes): Tffd700${Te6edf3eTb4b4b4.Te6edf3messageTffd700}Ta5d6ff"
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Wrap a [org.meshtastic.proto.TAKPacketV2] into the uncompressed v2 wire format: `[0xFF flag byte][raw protobuf]`.
* Used as a fallback when the zstd native lib isn't loaded.
*/
Tff7b72private Tff7b72fun Td2a8ffencodeUncompressedTb4b4b4(Te6edf3takPacketV2Tb4b4b4: Te6edf3orgTb4b4b4.Te6edf3meshtasticTb4b4b4.Te6edf3protoTb4b4b4.Te6edf3TAKPacketV2Tb4b4b4)Tb4b4b4: Te6edf3ByteArray Tb4b4b4{
Tff7b72val Te6edf3protoBytes Tff7b72= Te6edf3orgTb4b4b4.Te6edf3meshtasticTb4b4b4.Te6edf3protoTb4b4b4.Te6edf3TAKPacketV2Tb4b4b4.Te6edf3ADAPTERTb4b4b4.Te6edf3encodeTb4b4b4(Te6edf3takPacketV2Tb4b4b4)
Tff7b72val Te6edf3out Tff7b72= Te6edf3ByteArrayTb4b4b4(T79c0ff1 Tff7b72+ Te6edf3protoBytesTb4b4b4.Te6edf3sizeTb4b4b4)
Tff7b72outTff7b72[T79c0ff0Tff7b72] Tff7b72= Te6edf3TakV2CompressorTb4b4b4.Te6edf3DICT_ID_UNCOMPRESSEDTb4b4b4.Te6edf3toByteTb4b4b4(Tb4b4b4)
Te6edf3protoBytesTb4b4b4.Te6edf3copyIntoTb4b4b4(Tff7b72outTb4b4b4, T79c0ff1Tb4b4b4)
Tff7b72return Tff7b72out
Tb4b4b4}
T8b949e// ββ Receive: mesh β TAK client ββββββββββββββββββββββββββββββββββββββββββ
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffhandleMeshPacketTb4b4b4(Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3payload Tff7b72= Te6edf3packetTb4b4b4.Te6edf3decodedTff7b72?.Te6edf3payload Tff7b72?: Tff7b72return
Tff7b72when Tb4b4b4(Te6edf3packetTb4b4b4.Te6edf3decodedTff7b72?.Te6edf3portnumTb4b4b4) Tb4b4b4{
Te6edf3PortNumTb4b4b4.Te6edf3ATAK_PLUGIN_V2 Tff7b72-Tff7b72> Te6edf3handleV2PacketTb4b4b4(Te6edf3payloadTb4b4b4.Te6edf3toByteArrayTb4b4b4(Tb4b4b4)Tb4b4b4)
Te6edf3PortNumTb4b4b4.Te6edf3ATAK_PLUGIN Tff7b72-Tff7b72> Te6edf3handleV1PacketTb4b4b4(Te6edf3payloadTb4b4b4)
Tff7b72else Tff7b72-Tff7b72> Tff7b72return
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffhandleV2PacketTb4b4b4(Te6edf3wirePayloadTb4b4b4: Te6edf3ByteArrayTb4b4b4) Tb4b4b4{
Tff7b72try Tb4b4b4{
T8b949e// Decompress to CoT XML via the SDK's CotXmlBuilder, which handles
T8b949e// ALL typed payloads (DrawnShape, Marker, Route, etc.) and preserves
T8b949e// shape detail elements (vertices, colors, stroke weight) that the
T8b949e// app's own CoTXmlParser would strip. Forward the SDK-generated XML
T8b949e// directly to TAK clients without re-parsing.
Tff7b72val Te6edf3rawXml Tff7b72= Te6edf3TakV2CompressorTb4b4b4.Te6edf3decompressToXmlTb4b4b4(Te6edf3wirePayloadTb4b4b4)
T8b949e// Normalize for the TAK TCP stream β drop the <?xml ...?> prologue
T8b949e// and collapse inter-tag whitespace so ATAK's streaming parser sees
T8b949e// bare <event>...</event> on a single line. Centralized in the SDK.
Tff7b72val Te6edf3xml Tff7b72= Te6edf3CotSanitizerTb4b4b4.Te6edf3normalizeCotXmlTb4b4b4(Te6edf3rawXmlTb4b4b4)
T8b949e// Drop exact duplicates the mesh delivered more than once (multi-path relay or
T8b949e// retransmit) so ATAK doesn't surface doubled chat / TAK-Talk messages. Genuine
T8b949e// updates (new PLI position, moved marker, β¦) differ in content and pass through.
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3deliveryDedupTb4b4b4.Te6edf3admitTb4b4b4(Te6edf3xmlTb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffDropped duplicate CoT from mesh (already delivered within dedup window)Ta5d6ff" Tb4b4b4}
Tff7b72return
Tb4b4b4}
T8b949e// Logger.d { "RAW CoT IN (mesh): $xml" }
T8b949e// Routes: ATAK ignores b-m-r CoT events over TCP streaming.
T8b949e// Convert to a KML data package and write to ATAK's auto-import dir.
Tff7b72if Tb4b4b4(Te6edf3xmlTb4b4b4.Te6edf3containsTb4b4b4(Ta5d6ff"""Ta5d6fftype=Ta5d6ff"Ta5d6ffb-m-rTa5d6ff"""Ta5d6ff"Ta5d6ff)) {
try {
val pkg = RouteDataPackageGenerator.generateDataPackage(xml)
if (pkg != null) {
val (fileName, zipBytes) = pkg
AtakFileWriter.writeToImportDir(fileName, zipBytes)
} else {
Logger.w { Ta5d6ff"Te6edf3Route Te6edf3data Tff7b72package T7ee787generation Te6edf3failed β Te6edf3not Te6edf3enough Te6edf3waypoints?Ta5d6ff"Ta5d6ff }
}
} catch (e2: Exception) {
Logger.w(e2) { Ta5d6ff"Te6edf3Route Te6edf3data Tff7b72package T7ee787write Te6edf3failedTb4b4b4: $Tb4b4b4{Te6edf3e2Tb4b4b4.Te6edf3messageTb4b4b4}Ta5d6ff"Ta5d6ff }
}
}
takServerManager.broadcastRawXml(xml)
Logger.d { Ta5d6ff"Te6edf3V2 β Te6edf3TAK Te6edf3clients Tb4b4b4(Te6edf3raw Te6edf3XMLTb4b4b4)Ta5d6ff"Ta5d6ff }
} catch (e: Exception) {
Logger.w(e) { Ta5d6ff"Te6edf3Failed Te6edf3to Te6edf3handle Te6edf3V2 Te6edf3packetTb4b4b4: $Tb4b4b4{Te6edf3eTb4b4b4.Te6edf3messageTb4b4b4}Ta5d6ff"Ta5d6ff }
}
}
/**
* v1 receive path (firmware <= 2.7.x): decode bare protobuf [TAKPacket] (no compression) from port 72 (ATAK_PLUGIN)
* and convert to CoT for forwarding to attached TAK clients. Kept indefinitely so users on stable 2.7.x firmware
* retain PLI + GeoChat interop; new typed payloads (shapes, markers, routes, etc.) still require a v2-capable radio
* (firmware >= 2.8.0).
*/
private suspend fun handleV1Packet(payload: okio.ByteString) {
try {
val takPacket = TAKPacket.ADAPTER.decode(payload)
val cotMessage = convertV1ToCoT(takPacket) ?: return
takServerManager.broadcast(cotMessage)
Logger.d { Ta5d6ff"Te6edf3V1 β Te6edf3TAK Te6edf3clientsTb4b4b4: $Tb4b4b4{Te6edf3cotMessageTb4b4b4.Te6edf3typeTb4b4b4}Ta5d6ff"Ta5d6ff }
} catch (e: Exception) {
Logger.w(e) { Ta5d6ff"Te6edf3Failed Te6edf3to Te6edf3handle Te6edf3V1 Te6edf3packetTb4b4b4: $Tb4b4b4{Te6edf3eTb4b4b4.Te6edf3messageTb4b4b4}Ta5d6ff"Ta5d6ff }
}
}
private fun convertV1ToCoT(takPacket: TAKPacket): CoTMessage? {
val callsign = takPacket.contact?.callsign ?: Ta5d6ff"Te6edf3UNKNOWNTa5d6ff"
Ta5d6ff val senderUid = takPacket.contact?.device_callsign ?: Ta5d6ff"Te6edf3unknownTa5d6ff"
Ta5d6ff val teamName = takPacket.group?.team?.toTakTeamName() ?: DEFAULT_TAK_TEAM_NAME
val roleName = takPacket.group?.role?.toTakRoleName() ?: DEFAULT_TAK_ROLE_NAME
val battery = takPacket.status?.battery ?: DEFAULT_TAK_BATTERY
val pli = takPacket.pli
if (pli != null) {
return CoTMessage.pli(
uid = senderUid,
callsign = callsign,
latitude = pli.latitude_i.toDouble() / TAK_COORDINATE_SCALE,
longitude = pli.longitude_i.toDouble() / TAK_COORDINATE_SCALE,
altitude = pli.altitude.toDouble(),
speed = pli.speed.toDouble(),
course = pli.course.toDouble(),
team = teamName,
role = roleName,
battery = battery,
staleMinutes = DEFAULT_TAK_STALE_MINUTES,
)
}
val chat = takPacket.chat
if (chat != null) {
val timeNow = Clock.System.now()
// Include chatroom in UID so ATAK routes DMs correctly β the UID format
// Ta5d6ff"Te6edf3GeoChatTb4b4b4.Tff7b72<Te6edf3senderUidTff7b72>Tb4b4b4.Tff7b72<Te6edf3chatroomTff7b72>Tb4b4b4.Tff7b72<Te6edf3msgIdTff7b72>Ta5d6ff"Ta5d6ff is what ATAK uses to determine routing.
// Hardcoding Ta5d6ff"Te6edf3All Te6edf3Chat Te6edf3RoomsTa5d6ff"Ta5d6ff here loses DM routing from legacy v1 nodes.
val chatroom = chat.to ?: Ta5d6ff"Te6edf3All Te6edf3Chat Te6edf3RoomsTa5d6ff"
Ta5d6ff val msgId = Random.Default.nextInt().toString(TAK_HEX_RADIX)
return CoTMessage(
uid = Ta5d6ff"Te6edf3GeoChatTb4b4b4.Te6edf3$senderUidTb4b4b4.Te6edf3$chatroomTb4b4b4.Te6edf3$msgIdTa5d6ff"Ta5d6ff,
type = Ta5d6ff"Te6edf3bTff7b72-Te6edf3tTff7b72-Te6edf3fTa5d6ff"Ta5d6ff,
how = Ta5d6ff"Te6edf3hTff7b72-Te6edf3gTff7b72-Te6edf3iTff7b72-Te6edf3gTff7b72-Te6edf3oTa5d6ff"Ta5d6ff,
time = timeNow,
start = timeNow,
stale = timeNow + DEFAULT_TAK_STALE_MINUTES.minutes,
latitude = 0.0,
longitude = 0.0,
contact = CoTContact(callsign = callsign, endpoint = DEFAULT_TAK_ENDPOINT),
group = CoTGroup(name = teamName, role = roleName),
status = CoTStatus(battery = battery),
chat = CoTChat(chatroom = chatroom, senderCallsign = callsign, message = chat.message),
)
}
return null
}
companion object {
/**
* Minimum stale TTL (5 min) for static CoT types sent over mesh. iTAK uses 2-min stale for routes/shapes; over
* LoRa mesh with multi-hop relay, these arrive past stale and ATAK discards them. PLI and GeoChat are left
* untouched β their stale is meaningful.
*/
private val MIN_MESH_STALE_TTL = 15.minutes
private val STATIC_COT_PREFIXES = listOf(Ta5d6ff"Te6edf3bTff7b72-Te6edf3mTff7b72-Te6edf3rTa5d6ff"Ta5d6ff, Ta5d6ff"Te6edf3uTff7b72-Te6edf3dTff7b72-Ta5d6ff"Ta5d6ff, Ta5d6ff"Te6edf3bTff7b72-Te6edf3mTff7b72-Te6edf3pTff7b72-Ta5d6ff"Ta5d6ff)
private val EVENT_TYPE_RE = Regex(Ta5d6ff"Ta5d6ff"Ta5d6ff"Tff7b72<Te6edf3event\Te6edf3sTff7b72[Tff7b72^Tff7b72>Tff7b72]Tff7b72*\Te6edf3btypeTff7b72=Ta5d6ff"Ta5d6ff([^Ta5d6ff"Tff7b72]Tff7b72*Tb4b4b4)Ta5d6ff"""Ta5d6ff"Ta5d6ff)
// Matches the stale attribute ONLY within the <event> opening tag to avoid
// accidentally matching a stale=Ta5d6ff"Ta5d6ff...Ta5d6ff"Ta5d6ff on a <link> or other child element.
private val EVENT_TAG_RE = Regex(Ta5d6ff"""Tff7b72<Te6edf3event\Te6edf3bTff7b72[Tff7b72^Tff7b72>Tff7b72]Tff7b72*Tff7b72>Ta5d6ff"""Ta5d6ff)
private val STALE_ATTR_RE = Regex(Ta5d6ff"""\Te6edf3bstaleTff7b72=Ta5d6ff"Ta5d6ff([^Ta5d6ff"Tff7b72]Tff7b72*Tb4b4b4)Ta5d6ff"""Ta5d6ff"Ta5d6ff)
fun ensureMinimumStaleForMesh(xml: String): String {
val type = EVENT_TYPE_RE.find(xml)?.groupValues?.getOrNull(1) ?: return xml
if (STATIC_COT_PREFIXES.none { type.startsWith(it) }) return xml
// Search for stale only inside the <event> opening tag, not in child elements
val eventTagMatch = EVENT_TAG_RE.find(xml) ?: return xml
val eventTag = eventTagMatch.value
val staleInTag = STALE_ATTR_RE.find(eventTag) ?: return xml
val staleStr = staleInTag.groupValues[1]
val staleInstant =
try {
kotlin.time.Instant.parse(staleStr)
} catch (_: IllegalArgumentException) {
// Handle edge-case formats like missing Ta5d6ff"Ta5d6ffZTa5d6ff"
Ta5d6ff try {
val cleaned = staleStr.replace(Regex(Ta5d6ff"""\Tb4b4b4.\Te6edf3dTff7b72+Ta5d6ff"""Ta5d6ff), Ta5d6ff"Ta5d6ff"Ta5d6ff).replace(Ta5d6ff"Ta5d6ffZTa5d6ff"Ta5d6ff, Ta5d6ff"Ta5d6ff+00:00Ta5d6ff"Ta5d6ff)
kotlin.time.Instant.parse(cleaned)
} catch (_: IllegalArgumentException) {
return xml
}
}
val now = Clock.System.now()
val remaining = staleInstant - now
if (remaining >= MIN_MESH_STALE_TTL) return xml
val newStale = now + MIN_MESH_STALE_TTL
val newStaleStr = newStale.toString().replace(Regex(Ta5d6ff"""\Tb4b4b4.\Te6edf3dTff7b72+Ta5d6ff"""Ta5d6ff), Ta5d6ff"Ta5d6ff"Ta5d6ff) // strip fractional seconds
Logger.i {
Ta5d6ff"Ta5d6ffExtended stale for Tffd700$Te6edf3typeTa5d6ff: Tffd700$Te6edf3staleStrTa5d6ff β Tffd700$Te6edf3newStaleStrTa5d6ff Ta5d6ff"Ta5d6ff +
Ta5d6ff"Ta5d6ff(was Tffd700${Te6edf3remainingTb4b4b4.Te6edf3inWholeSecondsTffd700}Ta5d6ffs remaining, now Tffd700${Te6edf3MIN_MESH_STALE_TTLTb4b4b4.Te6edf3inWholeSecondsTffd700}Ta5d6ffs)Ta5d6ff"
Ta5d6ff }
// Replace the stale value only within the event tag, then splice the patched tag back
val newEventTag = eventTag.replaceRange(staleInTag.range, Ta5d6ff"""Te6edf3staleTff7b72=Ta5d6ff"Tffd700$Te6edf3newStaleStrTa5d6ff"Ta5d6ff"""Ta5d6ff)
return xml.replaceRange(eventTagMatch.range, newEventTag)
}
/**
* Strip non-essential CoT detail before mesh compression to save wire bytes. Delegates (via the [CotSanitizer]
* expect/actual seam) to the SDK's CotMeshSanitizer so the strip rules live in ONE golden-tested place shared
* by every consumer (Android, Apple, β¦) and can't drift between sides. Drift here once silently stripped
* TAK-Talk `<voice>` / `<marti>` and broke the feature end-to-end β guarded by the strip-preservation test in
* TAKMeshIntegrationTest and by CotMeshSanitizerTest in the SDK.
*/
fun stripNonEssentialElements(xml: String): String = CotSanitizer.stripNonEssentialForMesh(xml)
}
}
Served by rngit 1.5.0 - Generated in 0.09s